从简单案例开始 WebSocket 协议在握手阶段使用的是 HTTP 协议,握手成功后会升级(Upgrade)为双向通道,数据传输不再使用普通的字节流,而是使用名为 Frame(数据帧)的基本单位。Netty 贴心地为我们提供了很多开箱即用的编解码器,使得我们只需要几行配置就能搭建好它。以下是完整的服务端代码实现:
主引导类的编写 WebSocket 的核心在于它的 Pipeline(责任链)流水线配置。我们需要添加 HTTP 编解码器、聚合器,以及 Netty 自带的 WebSocket 协议服务器协议处理器。
MyWebSocketServer:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 import io.netty.bootstrap.ServerBootstrap;import io.netty.channel.ChannelFuture;import io.netty.channel.ChannelInitializer;import io.netty.channel.ChannelPipeline;import io.netty.channel.EventLoopGroup;import io.netty.channel.nio.NioEventLoopGroup;import io.netty.channel.socket.SocketChannel;import io.netty.channel.socket.nio.NioServerSocketChannel;import io.netty.handler.codec.http.HttpObjectAggregator;import io.netty.handler.codec.http.HttpServerCodec;import io.netty.handler.codec.http.websocketx.WebSocketServerProtocolHandler;import io.netty.handler.stream.ChunkedWriteHandler;public class MyWebSocketServer { public static void main (String[] args) { EventLoopGroup bossGroup = new NioEventLoopGroup (1 ); EventLoopGroup workerGroup = new NioEventLoopGroup (); try { ServerBootstrap serverBootstrap = new ServerBootstrap (); serverBootstrap.group(bossGroup, workerGroup) .channel(NioServerSocketChannel.class) .childHandler(new ChannelInitializer <SocketChannel>() { @Override protected void initChannel (SocketChannel ch) throws Exception { ChannelPipeline pipeline = ch.pipeline(); pipeline.addLast(new HttpServerCodec ()); pipeline.addLast(new ChunkedWriteHandler ()); pipeline.addLast(new HttpObjectAggregator (8192 )); pipeline.addLast(new WebSocketServerProtocolHandler ("/hello" )); pipeline.addLast(new MyWebSocketHandler ()); } }); System.out.println("WebSocket 服务端正在启动..." ); ChannelFuture channelFuture = serverBootstrap.bind(8888 ).sync(); System.out.println("WebSocket 服务端启动成功,监听端口:" + 8888 ); channelFuture.channel().closeFuture().sync(); } catch (InterruptedException e) { throw new RuntimeException (e); } finally { bossGroup.shutdownGracefully(); workerGroup.shutdownGracefully(); } } }
业务处理器的编写 MyWebSocketHandler:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 import io.netty.channel.ChannelHandlerContext;import io.netty.channel.SimpleChannelInboundHandler;import io.netty.handler.codec.http.websocketx.TextWebSocketFrame;import java.time.LocalDateTime;public class MyWebSocketHandler extends SimpleChannelInboundHandler <TextWebSocketFrame> { @Override protected void channelRead0 (ChannelHandlerContext ctx, TextWebSocketFrame msg) throws Exception { String text = msg.text(); System.out.println("【WebSocket服务端】收到来自客户端[" + ctx.channel().remoteAddress() + "]的消息: " + text); String responseText = "【服务器回复】" + LocalDateTime.now() + " -> 我已收到你的消息: " + text; ctx.channel().writeAndFlush(new TextWebSocketFrame (responseText)); } @Override public void channelActive (ChannelHandlerContext ctx) throws Exception { System.out.println("【WebSocket服务端】客户端连接成功: " + ctx.channel().remoteAddress()); } @Override public void channelInactive (ChannelHandlerContext ctx) throws Exception { System.out.println("【WebSocket服务端】客户端断开连接: " + ctx.channel().remoteAddress()); } @Override public void exceptionCaught (ChannelHandlerContext ctx, Throwable cause) throws Exception { System.err.println("【WebSocket服务端】发生异常" + cause); ctx.close(); } }
浏览器客户端测试 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 <!DOCTYPE html > <html lang ="zh-CN" > <head > <meta charset ="UTF-8" > <title > Netty WebSocket 测试客户端</title > <style > body { font-family : 'Segoe UI' , Tahoma, Geneva, Verdana, sans-serif; background-color : #f4f7f6 ; margin : 0 ; padding : 20px ; display : flex; justify-content : center; } .container { width : 100% ; max-width : 600px ; background : #ffffff ; padding : 25px ; border-radius : 10px ; box-shadow : 0 4px 15px rgba (0 ,0 ,0 ,0.08 ); } h2 { margin-top : 0 ; color : #333 ; border-bottom : 2px solid #eaeaea ; padding-bottom : 10px ; } .status-container { display : flex; align-items : center; margin-bottom : 15px ; font-size : 14px ; } .status-dot { width : 10px ; height : 10px ; border-radius : 50% ; background-color : #ff4d4f ; margin-right : 8px ; transition : background-color 0.3s ; } .status-dot .connected { background-color : #52c41a ; } #messageLog { height : 300px ; border : 1px solid #d9d9d9 ; background-color : #fafafa ; padding : 10px ; border-radius : 5px ; overflow-y : auto; margin-bottom : 15px ; font-family : 'Courier New' , Courier, monospace; font-size : 13px ; } .log-item { margin-bottom : 6px ; line-height : 1.4 ; } .log-inbound { color : #1890ff ; } .log-outbound { color : #52c41a ; } .log-system { color : #8c8c8c ; font-style : italic; } .input-group { display : flex; gap : 10px ; } #inputMsg { flex : 1 ; padding : 10px ; border : 1px solid #d9d9d9 ; border-radius : 4px ; outline : none; font-size : 14px ; } #inputMsg :focus { border-color : #40a9ff ; } button { padding : 0 20px ; background-color : #1890ff ; color : #fff ; border : none; border-radius : 4px ; cursor : pointer; font-size : 14px ; transition : background 0.2s ; } button :hover { background-color : #40a9ff ; } button :disabled { background-color : #f5f5f5 ; color : #rgba (0 , 0 , 0 , 0.25 ); border-color : #d9d9d9 ; cursor : not-allowed; } </style > </head > <body > <div class ="container" > <h2 > WebSocket 简易客户端</h2 > <div class ="status-container" > <div id ="statusDot" class ="status-dot" > </div > <span id ="statusText" > 正在尝试连接服务器...</span > </div > <div id ="messageLog" > </div > <div class ="input-group" > <input type ="text" id ="inputMsg" placeholder ="请输入要发送给 Netty 的消息..." disabled autocomplete ="off" > <button id ="sendBtn" onclick ="sendMessage()" disabled > 发送</button > </div > </div > <script > const serverUrl = "ws://127.0.0.1:8888/hello" ; let socket; const statusDot = document .getElementById ("statusDot" ); const statusText = document .getElementById ("statusText" ); const messageLog = document .getElementById ("messageLog" ); const inputMsg = document .getElementById ("inputMsg" ); const sendBtn = document .getElementById ("sendBtn" ); function initWebSocket ( ) { appendLog ("正在连接到: " + serverUrl, "system" ); socket = new WebSocket (serverUrl); socket.onopen = function (event ) { statusDot.classList .add ("connected" ); statusText.innerText = "已成功连接至 Netty 服务端" ; inputMsg.disabled = false ; sendBtn.disabled = false ; appendLog ("连接建立成功!可以开始发送消息了。" , "system" ); }; socket.onmessage = function (event ) { appendLog (event.data , "inbound" ); }; socket.onclose = function (event ) { statusDot.classList .remove ("connected" ); statusText.innerText = "与服务端断开连接" ; inputMsg.disabled = true ; sendBtn.disabled = true ; appendLog ("连接已断开。若要重新连接,请刷新页面。" , "system" ); }; socket.onerror = function (error ) { appendLog ("WebSocket 发生错误,请检查服务端是否正常启动。" , "system" ); }; } function sendMessage ( ) { const msg = inputMsg.value .trim (); if (!msg) return ; socket.send (msg); appendLog ("【我发送的】" + msg, "outbound" ); inputMsg.value = "" ; inputMsg.focus (); } function appendLog (text, type ) { const item = document .createElement ("div" ); item.className = "log-item log-" + type; const now = new Date ().toTimeString ().split (' ' )[0 ]; item.innerText = `[${now} ] ${text} ` ; messageLog.appendChild (item); messageLog.scrollTop = messageLog.scrollHeight ; } inputMsg.addEventListener ("keydown" , function (event ) { if (event.key === "Enter" ) { sendMessage (); } }); window .onload = initWebSocket; </script > </body > </html >
在 WebSocket 的建立过程中,101 状态码(101 Switching Protocols,切换协议)是整个长连接能够顺利诞生的“入场券”。 如果没有 101 状态码,WebSocket 根本无法启动。简单来说,客户端通过 HTTP 发起升级申请,服务端通过回传 101 状态码正式批准 “变身” 为 WebSocket 协议。为什么要这样设计呢?这是因为:
WebSocket 的设计目标是建立全双工、低延迟的长连接通道,但它在诞生之初面临一个现实问题:如何优雅地穿透互联网上各种复杂的防火墙和网关。如果直接在非标准端口上启动一套全新的 TCP 协议,很容易被各大路由节点当成恶意流量直接墙掉。为了兼容现有的 Web 基础设施,WebSocket 巧妙地选择了借鸡生蛋:
先通过标准的 HTTP 协议(80/443 端口)发起普通的握手请求。
握手一旦达成,立刻在原地将协议“升级切换”为 WebSocket 协议。
HTTP 规范中,专门定义了一个状态码来应对这种“中途换道”的场景,这就是 101 Switching Protocols。
我们可以通过客户端与 Netty 服务端在握手瞬间的实际抓包报文,来看 101 是如何起作用的:
阶段一:客户端发起升级申请(HTTP GET)。 当你在浏览器执行 new WebSocket(“ws://…”) 时,前端会发出一个特殊的 HTTP 请求,里面包含了“我要升级”的强烈信号:
1 2 3 4 5 6 7 GET /hello HTTP/1.1 Host: 127.0.0.1:8888 Upgrade: websocket <-- 【信号1】告诉服务器:我想把协议升级为 websocket Connection: Upgrade <-- 【信号2】告诉服务器:不要断开,我们要升级连接 Sec-WebSocket-Key: dGhlIHNhbXBsZSBub25jZQ== Origin: null Sec-WebSocket-Version: 13
阶段二:服务端同意升级,回传 101(HTTP 响应) 。Netty 服务端底层的 WebSocketServerProtocolHandler 接收到这个请求后,进行跨域和路径校验。如果通过,它会立刻原路返回一个 101 状态码的响应:
1 2 3 4 HTTP/1.1 101 Switching Protocols <-- 核心:老板同意了!协议现在正式切换! Upgrade: websocket <-- 确认升级为 websocket Connection: Upgrade <-- 确认连接保持并升级 Sec-WebSocket-Accept: s3pPLMBiTxaQ9kYGzzhZRbK+xOo=
一旦这个包含 101 的响应成功送达浏览器,原有的 HTTP 协议生命周期在这一瞬间彻底终结。在 Netty 层面,底层的流水线(Pipeline)会发生非常奇妙的变化:
老组件退场:专门处理 HTTP 解析的处理器(如 HttpServerCodec、HttpObjectAggregator)会自动功成身退。因为后续进来的不再是普通的 HTTP 报文,它们留在链条上只会报错。
新组件接管:WebSocketServerProtocolHandler 会把负责 WebSocket 帧解析(Frame Decoder/Encoder)的核心组件动态挂载到管道的最前线。
通道蜕变:此后,这个 TCP 通道上的内容不再有 GET、POST 或 Content-Type,取而代之的是由 Opcode(操作码)和 Payload(有效载荷)组成的 WebSocket 数据帧(Frame)。
加入心跳检测机制 下面我们对之前的 MyWebSocketServer 和业务 Handler 进行升级。规定:如果客户端连续 10 秒没有任何动静(读空闲),服务端就认为其掉线,主动关闭连接。
修改服务器引导类 我们需要在管道中挂载 IdleStateHandler。注意:它必须放在你的自定义业务处理器之前。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 pipeline.addLast(new IdleStateHandler (10 , 0 , 0 , TimeUnit.SECONDS)); pipeline.addLast(new WebSocketServerProtocolHandler ("/hello" )); pipeline.addLast(new MyWebSocketHandler ());
修改业务处理器 在 MyWebSocketHandler 增加重写实现方法 userEventTriggered,并适当修改消息响应:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 @Override public void userEventTriggered (ChannelHandlerContext ctx, Object evt) throws Exception { if (evt instanceof IdleStateEvent) { IdleStateEvent event = (IdleStateEvent) evt; if (event.state() == IdleState.READER_IDLE) { System.out.println("【心跳检测】客户端 [" + ctx.channel().remoteAddress() +"] 已经连续 10 秒没有任何读写行为,判定为僵尸连接!" ); ctx.channel().close(); } } else { super .userEventTriggered(ctx, evt); } } @Override protected void channelRead0 (ChannelHandlerContext ctx, TextWebSocketFrame msg) throws Exception { String text = msg.text(); System.out.println("【WebSocket服务端】收到来自客户端[" + ctx.channel().remoteAddress() + "]的消息: " + text); String responseText = "PING" .equals(text) ? "PONG" : "【服务器回复】" + LocalDateTime.now() + " -> 我已收到你的消息: " + text; ctx.channel().writeAndFlush(new TextWebSocketFrame (responseText)); }
升级 HTML 客户端 - 建立保活机制 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 let socket; let heartbeatIntervalId; // 用来存放定时器 ID function initWebSocket() { socket = new WebSocket("ws://127.0.0.1:8888/hello"); socket.onopen = function (event) { // ... 原有的 UI 解除置灰逻辑 ... appendLog("连接建立成功!", "system"); // 👉 开启客户端心跳保活定时器:每隔 5 秒向 Netty 发送一个特制的心跳字符串 startHeartbeat(); }; socket.onmessage = function (event) { // 过滤掉服务端返回的心跳回应,不让它污染用户的聊天日志 if (event.data === "PONG" || event.data.includes("心跳")) { return; } appendLog(event.data, "inbound"); }; socket.onclose = function (event) { // ... 原有的置灰逻辑 ... appendLog("连接已断开。", "system"); // 👉 连接断开了,必须清除定时器,防止死循环空发 stopHeartbeat(); }; } // 👉 启动心跳 function startHeartbeat() { heartbeatIntervalId = setInterval(() => { if (socket && socket.readyState === WebSocket.OPEN) { // 发送两端约定好的心跳帧文本,通常叫 "PING" socket.send("PING"); // 本地可以打个淡灰色的日志,证明客户端还在努力活着 console.log("[心跳] 向服务端发送 PING 保活..."); } }, 5000); // 5秒发一次,远小于服务端的 10 秒限制,安全 } // 👉 停止心跳 function stopHeartbeat() { if (heartbeatIntervalId) { clearInterval(heartbeatIntervalId); } }
支持连接鉴权 在实际生产环境中,连接鉴权(Authentication)是长连接系统的第一道安全防线。如果允许匿名用户随意建立 WebSocket 连接,服务器的内存和网络带宽很容易被恶意攻击者榨干。要实现像 ws://127.0.0.1:8888/hello?token=abcdefg&userId=1001 这样带参数的握手并在握手阶段进行安全校验,我们需要利用 Netty 责任链的特殊生命周期。
核心原理:利用 HTTP 阶段拦截鉴权 前面我们说到,WebSocket 在正式建立前,必须经历一个 HTTP 101 升级握手 阶段。这意味着:带参数的 URL 请求,本质上是一个标准的 HTTP GET 请求。
在 Netty 中,WebSocketServerProtocolHandler 一旦检测到匹配的路径(如 /hello),就会自动处理 101 升级。因此,我们必须在升级发生之前,拦截这个 HTTP 请求,解析出里面的 token 和 userId 并进行合法性校验。 如果校验失败,直接返回 HTTP 401 Unauthorized 并斩断连接,根本不给它升级成 WebSocket 的机会。
具体的实现 编写鉴权处理器 我们编写一个专门的 HTTP 入站处理器 AuthHandler。它继承自 ChannelInboundHandlerAdapter(处理普通的 HTTP 对象)。
在实际的分布式大型项目中,握手阶段的 token 通常建议使用 JWT(JSON Web Token) 。JWT 中本身就加密包裹了 userId、过期时间以及权限标识。AuthHandler 拦截到请求后,直接在内存中用秘钥解密 JWT。这样不需要查数据库或 Redis,就能在 0.1 毫秒内完成鉴权,极大减轻了长连接握手风暴(如服务器重启时百万客户端同时重连)对后台核心数据库的冲击!
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 import io.netty.channel.ChannelFutureListener;import io.netty.channel.ChannelHandlerContext;import io.netty.channel.ChannelInboundHandlerAdapter;import io.netty.handler.codec.http.*;import io.netty.util.AttributeKey;import lombok.extern.slf4j.Slf4j;import java.util.List;import java.util.Map;@Slf4j public class AuthHandler extends ChannelInboundHandlerAdapter { public static final AttributeKey<String> USER_ID_KEY = AttributeKey.valueOf("userId" ); @Override public void channelRead (ChannelHandlerContext ctx, Object msg) throws Exception { if (msg instanceof FullHttpRequest) { FullHttpRequest request = (FullHttpRequest) msg; String uri = request.uri(); if (uri.startsWith("/hello" )) { QueryStringDecoder decoder = new QueryStringDecoder (uri); Map<String, List<String>> parameters = decoder.parameters(); List<String> tokenList = parameters.get("token" ); List<String> userIdList = parameters.get("userId" ); String token = (tokenList != null && !tokenList.isEmpty()) ? tokenList.get(0 ) : null ; String userId = (userIdList != null && !userIdList.isEmpty()) ? userIdList.get(0 ) : null ; log.info("【鉴权层】收到握手申请,userId: {}, token: {}" , userId, token); if (!validateToken(userId, token)) { log.warn("【鉴权层】拒绝连接!用户 [{}] 鉴权失败。" , userId); sendUnauthorizedAndClose(ctx); return ; } ctx.channel().attr(USER_ID_KEY).set(userId); request.setUri(decoder.path()); } } super .channelRead(ctx, msg); } private boolean validateToken (String userId, String token) { if (userId == null || token == null ) { return false ; } return "abcdefg" .equals(token); } private void sendUnauthorizedAndClose (ChannelHandlerContext ctx) { FullHttpResponse response = new DefaultFullHttpResponse ( HttpVersion.HTTP_1_1, HttpResponseStatus.UNAUTHORIZED ); response.headers().set(HttpHeaderNames.CONTENT_TYPE, "text/plain; charset=UTF-8" ); HttpUtil.setContentLength(response, 0 ); ctx.writeAndFlush(response).addListener(ChannelFutureListener.CLOSE); } }
修改服务端流水线 在 MyWebSocketServer 中,我们将 AuthHandler 塞进 Pipeline。它的位置必须放在 HTTP 编解码和聚合器之后,但在 WebSocketServerProtocolHandler 之前!
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 @Override protected void initChannel (SocketChannel ch) throws Exception { ChannelPipeline pipeline = ch.pipeline(); pipeline.addLast(new HttpServerCodec ()); pipeline.addLast(new ChunkedWriteHandler ()); pipeline.addLast(new HttpObjectAggregator (8192 )); pipeline.addLast(new IdleStateHandler (10 , 0 , 0 , TimeUnit.SECONDS)); pipeline.addLast(new AuthHandler ()); pipeline.addLast(new WebSocketServerProtocolHandler ("/hello" )); pipeline.addLast(new MyWebSocketHandler ()); }
在业务层获取用户 当连接成功升级后,我们在聊天业务 MyWebSocketHandler 里就不需要每次都带 token 了。可以直接从当前 Channel 里取出刚才在 AuthHandler 里绑定的 userId:
1 2 3 4 5 6 7 8 9 10 11 @Override protected void channelRead0 (ChannelHandlerContext ctx, TextWebSocketFrame msg) throws Exception { String userId = ctx.channel().attr(AuthHandler.USER_ID_KEY).get(); String text = msg.text(); System.out.println("【WebSocket服务端】收到用户[" + userId + "]的消息: " + text); String responseText = "PING" .equals(text) ? "PONG" : "【服务器回复】" + LocalDateTime.now() + " -> [用户 " + userId + " ] 我已收到你的消息: " + text; ctx.channel().writeAndFlush(new TextWebSocketFrame (responseText)); }
升级 HTML 客户端 修改 HTML 里的 serverUrl,带上参数去连接。你可以分别尝试正确参数和错误参数:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 // 实验 1:使用正确合法的 token const serverUrl = "ws://127.0.0.1:8888/hello?token=abcdefg&userId=1001"; // 实验 2:修改 token 模拟黑客攻击(将其改为 "wrong_token") // const serverUrl = "ws://127.0.0.1:8888/hello?token=wrong_token&userId=1001"; // 监听连接关闭事件 socket.onclose = function (event) { statusDot.classList.remove("connected"); // 默认的系统提示 let reasonText = "与服务端断开连接"; /** * 利用 event.code 状态码来推断具体的错误原因! * 虽然 error 里没写,但关闭码里写得清清楚楚 */ if (event.code === 1006) { // 1006 代表异常关闭(非正常四次挥手断开) reasonText = "连接失败:服务端拒绝了握手请求(可能未启动、路径错误、或鉴权失败)"; appendLog(`[错误代码: ${event.code}] 握手被拒绝,请检查 token 是否正确、服务端是否存活。`, "system"); } else { appendLog(`[连接关闭] 关闭代码: ${event.code}, 原因: ${event.reason || '未指定'}`, "system"); } statusText.innerText = reasonText; inputMsg.disabled = true; sendBtn.disabled = true; }; // 监听发生错误事件 socket.onerror = function (error) { // 直接打印 error 拿不到具体字符串,但通过控制台对象展开可以看到其类型 console.error("【WebSocket 底层错误触发】:", error); appendLog("【网络提示】浏览器底层检测到 WebSocket 异常,连接即将关闭...", "system"); };
实验1:
实验2:
群聊与私聊的实现 在原生的 NIO 中,为了管所有客户端的 SocketChannel,我们通常得自己写一个 ConcurrentHashMap 来维护用户 ID 和 Channel 的映射关系。到了 Netty 中,框架为我们提供了一个极其强大的容器 —— ChannelGroup。它可以帮我们把成千上万的 Channel 组织成一个聚合组。不仅能自动管理 Channel 的生命周期(当客户端断开连接时,Netty 会自动将其从 ChannelGroup 中移除,不需要你手动 remove),还提供了优雅的、批量的高性能异步广播机制。下面我们结合前面的 AuthHandler(已将 userId 绑定到 Channel 属性中),来实现一个支持精准单聊(私聊)和全员广播(群聊)的案例。
编写全网连接管理器 在实际开发中,我们通常会抽象出一个单例的 SessionManager,用来同时兼顾 “全员群聊(ChannelGroup)” 和 “精准私聊(ConcurrentHashMap)”。
MySessionManager:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 import io.netty.channel.Channel;import io.netty.channel.group.ChannelGroup;import io.netty.channel.group.DefaultChannelGroup;import io.netty.handler.codec.http.websocketx.TextWebSocketFrame;import io.netty.util.concurrent.GlobalEventExecutor;import lombok.extern.slf4j.Slf4j;import java.util.concurrent.ConcurrentHashMap;@Slf4j public class MySessionManager { private MySessionManager () {} private static final MySessionManager INSTANCE = new MySessionManager (); public static MySessionManager getInstance () { return INSTANCE; } private final ChannelGroup globalGroup = new DefaultChannelGroup (GlobalEventExecutor.INSTANCE); private final ConcurrentHashMap<String, Channel> userChannelMap = new ConcurrentHashMap <>(); public void register (String userId, Channel channel) { globalGroup.add(channel); userChannelMap.put(userId, channel); log.info("【会话管理】用户 [{}] 登录上线,当前在线人数: {}" , userId, userChannelMap.size()); } public void unregister (Channel channel) { String userId = channel.attr(AuthHandler.USER_ID_KEY).get(); if (userId != null ) { userChannelMap.remove(userId); } globalGroup.remove(channel); log.info("【会话管理】用户 [{}] 离线,当前在线人数: {}" , userId, userChannelMap.size()); } public void broadcast (String message) { globalGroup.writeAndFlush(new TextWebSocketFrame ("【系统广播】" + message)); } public void sendToUser (String toUserId, String message) { Channel targetChannel = userChannelMap.get(toUserId); if (targetChannel != null && targetChannel.isActive()) { targetChannel.writeAndFlush(new TextWebSocketFrame (message)); } else { log.warn("【会话管理】发送失败:用户 [{}] 当前不在线或连接已失效。" , toUserId); } } }
业务处理器升级 - 约定聊天协议以及用户上线离线的绑定 为了让服务端知道这条消息是“群发”还是“私聊给某人”,我们规定一个最简单的文本协议格式(通常生产环境用 JSON,这里用文本切分演示):
群聊格式:ALL:大家好呀!
私聊格式:TO:1002:嘿,哥们,你在吗?
修改:MyWebSocketHandler
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 public class MyWebSocketHandler extends SimpleChannelInboundHandler <TextWebSocketFrame> { @Override public void userEventTriggered (ChannelHandlerContext ctx, Object evt) throws Exception { if (evt instanceof IdleStateEvent) { IdleStateEvent event = (IdleStateEvent) evt; if (event.state() == IdleState.READER_IDLE) { System.out.println("【心跳检测】客户端 [" + ctx.channel().remoteAddress() + "] 已经连续 10 秒没有任何读写行为,判定为僵尸连接!" ); ctx.channel().close(); } } else if (evt instanceof WebSocketServerProtocolHandler.HandshakeComplete) { String userId = ctx.channel().attr(AuthHandler.USER_ID_KEY).get(); if (userId != null ) { MySessionManager.getInstance().register(userId, ctx.channel()); MySessionManager.getInstance().broadcast("欢迎用户 [" + userId + "] 进入聊天室!" ); } } else { super .userEventTriggered(ctx, evt); } } @Override protected void channelRead0 (ChannelHandlerContext ctx, TextWebSocketFrame msg) throws Exception { String currentUserId = ctx.channel().attr(AuthHandler.USER_ID_KEY).get(); String rawText = msg.text(); if ("PING" .equals(rawText)) { ctx.channel().writeAndFlush(new TextWebSocketFrame ("PONG" )); return ; } System.out.println("【业务层】收到用户 [" + currentUserId +"] 的原始消息: " + rawText); if (rawText.startsWith("ALL:" )) { String content = rawText.substring(4 ); String broadcastMsg = String.format("【群聊】用户[%s]说: %s" , currentUserId, content); MySessionManager.getInstance().broadcast(broadcastMsg); } else if (rawText.startsWith("TO:" )) { String[] parts = rawText.split(":" , 3 ); if (parts.length == 3 ) { String toUserId = parts[1 ]; String content = parts[2 ]; String privateMsg = String.format("【私聊】用户[%s]悄悄对你说: %s" , currentUserId, content); MySessionManager.getInstance().sendToUser(toUserId, privateMsg); ctx.channel().writeAndFlush(new TextWebSocketFrame (String.format("【你对%s说】: %s" , toUserId, content))); } else { ctx.channel().writeAndFlush(new TextWebSocketFrame ("【系统提示】格式错误,私聊请使用: TO:用户ID:内容" )); } } else { ctx.channel().writeAndFlush(new TextWebSocketFrame ("【系统提示】不认识的指令。群发请输入 ALL:内容,私聊请输入 TO:用户ID:内容" )); } } @Override public void channelActive (ChannelHandlerContext ctx) throws Exception { System.out.println("【系统提示】底层通道已激活: " + ctx.channel().remoteAddress()); } @Override public void channelInactive (ChannelHandlerContext ctx) throws Exception { String userId = ctx.channel().attr(AuthHandler.USER_ID_KEY).get(); MySessionManager.getInstance().unregister(ctx.channel()); if (userId != null ) { MySessionManager.getInstance().broadcast("用户 [" + userId + "] 离开了聊天室。" ); } } @Override public void exceptionCaught (ChannelHandlerContext ctx, Throwable cause) throws Exception { System.err.println("【WebSocket服务端】发生异常" + cause); ctx.close(); } }
注意,上述 MySessionManager.getInstance().register 必须在 userEventTriggered 才可以(而不是 channelActive),这是网络 TCP 连接的建立过程以及Netty 底层的事件流向决定的。简单来说:
channelActive 代表的是 “物理通道(TCP)通了”;
而 userEventTriggered 里的 HandshakeComplete 代表的是 “业务协议(WebSocket)谈成了”。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 [步骤 1] 客户端发起 TCP 三次握手 │ [步骤 2] 🤝 TCP 握手成功!操作系统的 Socket 通道建立 │ └─► 🔥 触发 channelActive() (此时:通道里空空如也,客户端连 1 个字节的 HTTP 握手数据都还没发过来) │ [步骤 3] 📨 浏览器发送 HTTP 请求报文(GET /hello?userId=1001&token=...) │ [步骤 4] ⚙️ 数据流入 Pipeline:HttpServerCodec ──► AuthHandler (解析出 userId 并解析绑定) │ [步骤 5] 🎯 WebSocketServerProtocolHandler 进行 101 状态码协议升级响应 │ └─► 🔥 触发 userEventTriggered() 并且携带 HandshakeComplete 事件 (此时:WebSocket 长连接正式宣告诞生)
联调测试 现在你可以同时打开两个不同的浏览器窗口(或者一个是正常模式,一个是无痕模式),模拟两个不同的用户同时在线:
玩家 A 上线:
修改 HTML 的 URL 里的 userId 为 1001: const serverUrl = “ws://127.0.0.1:8888/hello?token=abcdefg&userId=1001”;
打开网页 A。系统日志会显示:[系统广播] 欢迎用户 [1001] 进入聊天室!
玩家 B 上线:
修改同一个 HTML 里的 userId 为 1002: const serverUrl = “ws://127.0.0.1:8888/hello?token=abcdefg&userId=1002”;
打开网页 B。此时网页 A 和网页 B 都会同时收到:[系统广播] 欢迎用户 [1002] 进入聊天室!(证明 ChannelGroup 广播成功!)
测试群聊:
在用户 1001 的输入框输入:ALL:大家好,我是1001!
结果:两个网页同时收到了这行群发。
测试精准私聊:
在用户 1001 的输入框输入:TO:1002:哥们,单独聊聊,今晚去吃烧烤吗?
观察结果:
用户 1002 的网页上:弹出了私聊字样:【私聊】用户[1001]悄悄对你说: 今晚去吃烧烤吗?
其他任何不相干用户的网页上:一片安静,什么都收不到。这就是精准定向推送。
用户上线:
广播消息:
单独私聊:
性能深度优化 在高并发的 WebSocket 推送系统中(如百万级在线的聊天室、直播弹幕、即时股票行情),内存暴涨和 OOM(内存溢出) 往往是除了线程阻塞之外最致命的杀手。Netty 为了榨干网络 IO 的性能,默认使用了池化的堆外内存(PooledByteBufAllocator)。堆外内存不受 JVM 垃圾回收(GC)的管辖,如果开发者在使用、跨线程传递过程中少释放了一次,就会造成底层的内核内存泄漏;如果多释放了一次,则会引发严重的 IllegalReferenceCountException: refCnt: 0 崩溃。针对以上高并发推送场景,我们将从 TextWebSocketFrame 的复用 和 跨线程派发 两个维度,深度剖析如何避开内存泄漏的坑。
跨线程池派发任务时的引用计数坑 在 Netty 中,所有的网络出入站对象(如 ByteBuf、TextWebSocketFrame 等)都实现了 ReferenceCounted(引用计数)接口。其核心铁律是:
对象的初始引用计数为 1。
只要引用计数变为 0,对象占用的底层堆外内存就会被立刻释放并回收到对象池。
谁消费,谁释放。 Netty 的 SimpleChannelInboundHandler 会在 channelRead0 结束后自动帮你把 msg 的引用计数减 1。
经典的跨线程泄漏死锁的例子:很多开发者在收到消息后,为了不阻塞 IO 线程(EventLoop),会把消息提交到自定义的业务线程池中处理:
1 2 3 4 5 6 7 8 9 10 11 @Override protected void channelRead0 (ChannelHandlerContext ctx, TextWebSocketFrame msg) { businessThreadPool.execute(() -> { String text = msg.text(); System.out.println(text); }); }
正确的跨线程传递姿势(保留引用):如果必须把包装了 ByteBuf 的对象(如 TextWebSocketFrame)传给其他线程,你必须显式调用 retain() 来增加引用计数,告诉 Netty:“我还要用,你先别回收!”。并在业务线程消费完毕后手动 release()。
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 @Override protected void channelRead0 (ChannelHandlerContext ctx, TextWebSocketFrame msg) { msg.retain(); businessThreadPool.execute(() -> { try { String text = msg.text(); } finally { msg.release(); } }); }
高并发下群发推送时 TextWebSocketFrame 极致复用 假设一个聊天室里有 1 万名在线用户,有人发了一条消息,服务端需要同时推送给这 1 万个人。
❌ 愚蠢的做法:循环 new 对象。这种做法在高并发下会产生大量的内存碎片,频繁触发池化分配器的锁竞争,甚至导致 OOM。
1 2 3 4 for (Channel ch : globalGroup) { ch.writeAndFlush(new TextWebSocketFrame ("大家好!" )); }
✅ 高性能复用方案:使用 ChannelGroup 或手动 retain()。Netty 的 ChannelGroup.writeAndFlush(object) 内部已经做过了极致的复用优化:它不会为每个通道单独复制一份数据,而是通过共享同一个底层 ByteBuf 内存,并对引用计数进行递增的方式实现的。如果你因为特殊的路由逻辑(比如按照用户标签推送),无法使用 ChannelGroup,必须自己写 for 循环群发,那么你必须采用 “单对象 + 手动 retain” 的零拷贝复用策略:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 public void broadcastToUsers (List<Channel> targetChannels, String message) { if (targetChannels == null || targetChannels.isEmpty()) return ; TextWebSocketFrame frame = new TextWebSocketFrame (message); int targetCount = targetChannels.size(); for (int i = 0 ; i < targetCount - 1 ; i++) { frame.retain(); } for (Channel ch : targetChannels) { if (ch.isActive()) { ch.writeAndFlush(frame); } else { frame.release(); } } }
监控和排查内存泄漏 即使代码写得再小心,高并发下依然可能存在链路遗漏导致的泄漏。Netty 提供了工业级的内存泄漏采样探测器。在生产环境或本地压力测试启动时,通过 JVM 参数开启最高级别的泄漏检测:
1 2 3 4 5 -Dio.netty.leakDetection.level=PARANOID
如果发生泄漏,控制台会弹出类似如下的恐怖警告:
1 2 3 4 5 LEAK: ByteBuf.release() was not called before it's garbage-collected. See https://netty.io/wiki/reference-counted-objects.html for more information. Recent access records: 4 #1: com.owlias.server.zdemo.AuthHandler.channelRead(AuthHandler.java:42) #2: io.netty.handler.codec.MessageToMessageDecoder.channelRead(...)
一旦看到这个日志,不要慌,顺着它提示的 #1、#2 堆栈,去找最后一次创建或调用 retain(),但后续没有安全闭环 release() 的代码位置,将其放进 try-finally 块中确保释放即可。
标题:
Java NIO - Netty 支持 WebSocket 详细案例以及深度优化方案